Skip to main content

media_pp\elements\sink/
frame_counter.rs

1use std::sync::{
2    Arc,
3    atomic::{AtomicUsize, Ordering},
4};
5
6use crate::pp_log::{PpLog, pp_info};
7
8use crate::{
9    buffer::MediaBuffer,
10    control::ControlMsg,
11    element::{Element, ElementType, Sink, element_pp_log},
12    error::Result,
13};
14
15/// Terminal sink that counts decoded frames (video or audio). Backed by
16/// an `Arc<AtomicUsize>` so the count can be read from outside the
17/// pipeline even when this sink ends up running on a `Queue` worker
18/// thread.
19pub struct FrameCounter {
20    pp_log: PpLog,
21    name: Arc<str>,
22    count: Arc<AtomicUsize>,
23}
24
25impl FrameCounter {
26    pub fn new(name: impl Into<String>) -> (Self, Arc<AtomicUsize>) {
27        let count = Arc::new(AtomicUsize::new(0));
28        let name: Arc<str> = name.into().into();
29        let pp_log = element_pp_log(ElementType::FrameCounter, &name, None);
30        pp_info!(pp_log: &pp_log, "created");
31        (
32            Self {
33                name,
34                pp_log,
35                count: count.clone(),
36            },
37            count,
38        )
39    }
40}
41
42impl Element for FrameCounter {
43    fn name(&self) -> Arc<str> {
44        self.name.clone()
45    }
46
47    fn element_type(&self) -> ElementType {
48        ElementType::FrameCounter
49    }
50
51    fn pp_log(&self) -> &PpLog {
52        &self.pp_log
53    }
54
55    fn pp_log_mut(&mut self) -> &mut PpLog {
56        &mut self.pp_log
57    }
58}
59
60impl Sink for FrameCounter {
61    fn consume(&mut self, buf: MediaBuffer) -> Result<()> {
62        if let MediaBuffer::Video(_) | MediaBuffer::Audio(_) = buf {
63            self.count.fetch_add(1, Ordering::Relaxed);
64        }
65        Ok(())
66    }
67
68    fn control(&mut self, _msg: ControlMsg) -> Result<()> {
69        // Terminal, nothing to flush or forward.
70        Ok(())
71    }
72}